Skip to main content
Version: Next

Azure OpenAI for big data

The Azure OpenAI service can be used to solve a large number of natural language tasks through chat, responses, and embedding APIs. To make it easier to scale your prompting workflows from a few examples to large datasets of examples, we have integrated the Azure OpenAI service with the distributed machine learning library SynapseML. This integration makes it easy to use the Apache Spark distributed computing framework to process millions of prompts with the OpenAI service. This tutorial shows how to apply large language models at a distributed scale using Azure OpenAI.

Prerequisites​

The key prerequisites for this quickstart include a working Azure OpenAI resource, and an Apache Spark cluster with SynapseML installed. We suggest creating a Synapse workspace, but an Azure Databricks, HDInsight, or Spark on Kubernetes, or even a python environment with the pyspark package will work.

  1. An Azure OpenAI resource – request access here before creating a resource
  2. Create a Synapse workspace
  3. Create a serverless Apache Spark pool

Import this guide as a notebook​

The next step is to add this code into your Spark cluster. You can either create a notebook in your Spark platform and copy the code into this notebook to run the demo. Or download the notebook and import it into Synapse Analytics

Fill in service information​

Next, edit the cell in the notebook to point to your service. In particular set the service_name, deployment_name, location, and key variables to match them to your OpenAI service:

from synapse.ml.core.platform import find_secret

# Fill in the following lines with your service information
# Learn more about selecting which embedding model to choose: https://openai.com/blog/new-and-improved-embedding-model
service_name = "synapseml-openai-3"
deployment_name = "gpt-5-mini"
deployment_name_embeddings = "text-embedding-ada-002"

key = find_secret(
secret_name="openai-api-key-3", keyvault="mmlspark-build-keys"
) # please replace this line with your key as a string

assert key is not None and service_name is not None

Create a dataset of prompts​

Next, create a dataframe consisting of a series of rows, with one prompt per row.

You can also load data directly from ADLS or other databases. For more information on loading and preparing Spark dataframes, see the Apache Spark data loading guide.

# spark session is assumed to be created in the environment already such as in Fabric notebooks with Spark environment
df = spark.createDataFrame(
[
("Hello my name is",),
("The best code is code thats",),
("SynapseML is ",),
]
).toDF("prompt")

More Usage Examples​

Generating Text Embeddings​

In addition to completing text, we can also embed text for use in downstream algorithms or vector retrieval architectures. Creating embeddings allows you to search and retrieve documents from large collections and can be used when prompt engineering isn't sufficient for the task.

For more information on using OpenAIEmbedding see our embedding guide.

from synapse.ml.services.openai import OpenAIEmbedding

embedding = (
OpenAIEmbedding()
.setSubscriptionKey(key)
.setDeploymentName(deployment_name_embeddings)
.setCustomServiceName(service_name)
.setTextCol("prompt")
.setErrorCol("error")
.setOutputCol("embeddings")
)

display(embedding.transform(df).show(truncate=False))

Chat Completion​

Models such as ChatGPT and GPT-4 are capable of understanding chats instead of single prompts. The OpenAIChatCompletion transformer exposes this functionality at scale.

from synapse.ml.services.openai import OpenAIChatCompletion
from pyspark.sql import Row
from pyspark.sql.types import *


def make_message(role, content):
return Row(role=role, content=content, name=role)


chat_df = spark.createDataFrame(
[
(
[
make_message(
"system", "You are an AI chatbot with red as your favorite color"
),
make_message("user", "Whats your favorite color"),
],
),
(
[
make_message("system", "You are very excited"),
make_message("user", "How are you today"),
],
),
]
).toDF("messages")


chat_completion = (
OpenAIChatCompletion()
.setSubscriptionKey(key)
.setDeploymentName(deployment_name)
.setCustomServiceName(service_name)
.setMessagesCol("messages")
.setOutputCol("chat_completions")
.setErrorCol("chat_completions_error")
)

display(
chat_completion.transform(chat_df)
.select("messages", "chat_completions.choices.message.content")
.show(truncate=False)
)

Multimodal Chat Completions​

Run this example with a vision-capable deployment. Each user message contains a text part and an image URL part. Define an explicit Spark schema so the nested image_url.url and optional detail fields are preserved when OpenAIChatCompletion builds the request.

from pyspark.sql import Row
from pyspark.sql.types import ArrayType, StringType, StructField, StructType
from synapse.ml.services.openai import OpenAIChatCompletion

image_url_schema = StructType(
[
StructField("url", StringType(), False),
StructField("detail", StringType(), True),
]
)
content_part_schema = StructType(
[
StructField("type", StringType(), False),
StructField("text", StringType(), True),
StructField("image_url", image_url_schema, True),
]
)
message_schema = StructType(
[
StructField("role", StringType(), False),
StructField("content", ArrayType(content_part_schema), False),
StructField("name", StringType(), True),
]
)
multimodal_schema = StructType(
[StructField("messages", ArrayType(message_schema), False)]
)

multimodal_chat_df = spark.createDataFrame(
[
Row(
messages=[
Row(
role="user",
content=[
Row(
type="text",
text="Describe this image in one sentence.",
image_url=None,
),
Row(
type="image_url",
text=None,
image_url=Row(
url="https://mmlspark.blob.core.windows.net/datasets/OCR/test2.png",
detail="low",
),
),
],
name=None,
)
]
)
],
schema=multimodal_schema,
)

multimodal_chat = (
OpenAIChatCompletion()
.setSubscriptionKey(key)
.setDeploymentName(deployment_name)
.setCustomServiceName(service_name)
.setMessagesCol("messages")
.setOutputCol("multimodal_chat_completions")
.setErrorCol("multimodal_chat_errors")
)

multimodal_chat.transform(multimodal_chat_df).select(
"messages",
"multimodal_chat_completions.choices.message.content",
"multimodal_chat_errors",
).show(truncate=False)

Chat Completion with Reasoning Models​

GPT-5.1 and GPT-5-mini reasoning deployments reject explicit sampling parameters when reasoning is enabled:

  • Leave temperature, top_p, and seed unset.
  • Use reasoning_effort and verbosity to control reasoning behavior.

The deployment, credentials, and endpoint can still be set globally using OpenAIDefaults.

from synapse.ml.services.openai import OpenAIChatCompletion
from synapse.ml.services.openai.OpenAIDefaults import OpenAIDefaults

# Set global defaults without unsupported sampling parameters
defaults = OpenAIDefaults()
defaults.set_deployment_name(deployment_name)
defaults.set_subscription_key(key)
defaults.set_URL(f"https://{service_name}.openai.azure.com/")

chat_completion = (
OpenAIChatCompletion()
.setMessagesCol("messages")
.setOutputCol("chat_completions")
.setErrorCol("chat_completions_error")
)

display(
chat_completion.transform(chat_df)
.select("messages", "chat_completions.choices.message.content")
.show(truncate=False)
)

Retired Completions API​

The OpenAICompletion transformer has been removed because the legacy Completions API is deprecated and retired. Use OpenAIChatCompletion, OpenAIPrompt with chat_completions or responses, or OpenAIResponses for text generation workloads.

Structured JSON output with json_schema​

Use the new responseFormat option json_schema to force the model to return JSON strictly matching a provided JSON Schema. This improves reliability for downstream parsing and analytics.

Key points:

  • Set responseFormat to a dict containing type="json_schema" and a nested json_schema object (JSON string form not supported).
  • Provide the schema as a Python dict / Scala Map plus a name and optional strict flag.
  • Bare json_schema string is rejected; supply the full dict form.
  • Allowed simple String values for responseFormat: text, json_object.
  • For json_schema minimal validation only checks the presence of the json_schema key; the nested schema is passed through unchanged.
  • The model output will be constrained to the schema (no extra properties when additionalProperties: false).

Summary table

TypeHow to setRequires nested schemaNotes
textString("text")NoRaw string output
json_objectString("json_object")NoModel attempts well‑formed JSON (not strictly validated)
json_schemaDict/Map onlyYesStrict; reject bare string or JSON string form

Below we request a single field answer as structured JSON.

from synapse.ml.services.openai import OpenAIChatCompletion
from pyspark.sql import Row

# Define the JSON Schema we want the model to satisfy
schema = {
"type": "object",
"properties": {"answer": {"type": "string"}},
"required": ["answer"],
"additionalProperties": False,
}

# Single user message requesting structured output
messages_df = spark.createDataFrame(
[
(
[
Row(
role="user",
content="What is the capital of France?",
)
],
)
]
).toDF("messages")

chat_structured = (
OpenAIChatCompletion()
.setSubscriptionKey(key)
.setDeploymentName(deployment_name)
.setCustomServiceName(service_name)
.setMessagesCol("messages")
.setResponseFormat(
{
"type": "json_schema",
"json_schema": {
"name": "answer_schema",
"strict": True,
"schema": schema,
},
}
)
.setOutputCol("chat_structured")
.setErrorCol("chat_structured_error")
)

display(
chat_structured.transform(messages_df)
.select("chat_structured.choices.message.content")
.show(truncate=False)
)
# The returned content should be a JSON object: {\"answer\": \"Paris\"}

Schema-only structured output​

In development builds with setResponseSchema, you can pass the inner JSON Schema directly instead of constructing the response-format envelope. The helper defaults to name="response_schema" and strict=True, and accepts an optional name and strict flag.

OpenAIChatCompletion, OpenAIResponses, and OpenAIPrompt all support this method. Each stage uses the correct envelope for its API. Existing setResponseFormat behavior is unchanged.

Desired outputConfiguration
Plain textsetResponseFormat("text") or setResponseFormat({"type": "text"})
An arbitrary JSON objectsetResponseFormat("json_object") or setResponseFormat({"type": "json_object"})
A supplied inner JSON SchemasetResponseSchema(schema)
An already wrapped or named response formatsetResponseFormat(response_format)

The text and JSON-object selectors do not acquire a schema envelope. Named partial formats and full Chat/Responses envelopes retain their supplied name and strictness. The legacy setResponseFormat(inner_schema) form remains supported without implicitly enabling strict mode; use setResponseSchema when you want the strict-by-default convenience API. Do not pass an already wrapped format to setResponseSchema. Pydantic conversion remains the caller's responsibility.

An unchanged static response format is serialized lazily once per stage instance and its JSON is reused for later rows. The bounded cache is replaced when the format changes and is not persisted or transferred to Spark executors; each executor-side stage instance builds its own cache. Python dictionaries and lists are copied into immutable JVM collections; Scala formats containing mutable sequences remain uncached so later changes are reflected in requests. Column-driven response formats remain row-dependent, and Responses verbosity can still vary by row while its static format is cached. The format is still included in every HTTP request; caching does not cache model responses or reduce the number of requests.

The helper does not modify the JSON Schema or infer missing fields. For strict output, supply the model's supported schema, including required and additionalProperties: False for object definitions. Pass a Python dictionary or Scala Map[String, Any], not a JSON string or Spark StructType. Dictionary keys must be strings. Scala BigInt and BigDecimal values are not supported by the existing parameter serializer and are rejected before configuration changes. Continue to handle errors, refusals, and incomplete responses.

from synapse.ml.services.openai import OpenAIPrompt, OpenAIResponses
from pyspark.sql import functions as F

# Reuse the inner schema, credentials, and one-row messages_df above.
country_df = spark.createDataFrame([("France",)], ["country"])
schema_e2e_cases = []


def configure_schema_stage(stage, options):
return (
stage.setSubscriptionKey(key)
.setDeploymentName(deployment_name)
.setCustomServiceName(service_name)
.setResponseSchema(schema, **options)
.setOutputCol("schema_output")
.setErrorCol("schema_error")
)


def check_schema_result(case, result, parsed):
rows = result.select("schema_error", parsed.alias("parsed")).collect()
assert len(rows) == 1, f"{case}: expected one row, got {len(rows)}"
error = rows[0]["schema_error"]
assert error is None, f"{case}: service error: {error}"
parsed_answer = rows[0]["parsed"]
assert (
parsed_answer is not None
), f"{case}: expected parsed answer struct, got {parsed_answer!r}"
answer = parsed_answer["answer"]
assert (
isinstance(answer, str) and answer.strip()
), f"{case}: expected nonempty answer string, got {answer!r}"
schema_e2e_cases.append(case)
print(f"Schema-only E2E passed: {case}")


for option_name, options in (
("defaults", {}),
("custom_non_strict", {"name": "answer_schema", "strict": False}),
):
chat_schema = configure_schema_stage(OpenAIChatCompletion(), options)
chat_result = chat_schema.setMessagesCol("messages").transform(messages_df)
chat_text = F.expr("try_element_at(schema_output.choices, 1).message.content")
check_schema_result(
f"chat/{option_name}", chat_result, F.from_json(chat_text, "answer STRING")
)

responses_schema = configure_schema_stage(OpenAIResponses(), options)
responses_result = (
responses_schema.setMessagesCol("messages")
.setApiVersion("2025-04-01-preview")
.transform(messages_df)
)
responses_text = F.expr(
"try_element_at(filter(flatten(transform("
"filter(schema_output.output, x -> x.content IS NOT NULL), x -> x.content)), "
"x -> x.type = 'output_text'), 1).text"
)
check_schema_result(
f"responses/{option_name}",
responses_result,
F.from_json(responses_text, "answer STRING"),
)

for api_type in ("chat_completions", "responses"):
prompt_schema = configure_schema_stage(OpenAIPrompt(), options)
prompt_result = (
prompt_schema.setApiType(api_type)
.setApiVersion("2025-04-01-preview")
.setPromptTemplate("What is the capital of {country}?")
.setPostProcessing("json")
.setPostProcessingOptions({"jsonSchema": "answer STRING"})
.transform(country_df)
)
check_schema_result(
f"prompt_{api_type}/{option_name}", prompt_result, F.col("schema_output")
)

assert len(schema_e2e_cases) == 8

Scala uses the same schema-only API. Use scala.collection.immutable.ListMap for each object whose key order matters; an ordinary Map does not guarantee insertion order:

import com.microsoft.azure.synapse.ml.services.openai.OpenAIChatCompletion
import scala.collection.immutable.ListMap

val schema: Map[String, Any] = ListMap(
"type" -> "object",
"properties" -> ListMap("answer" -> Map("type" -> "string")),
"required" -> Seq("answer"),
"additionalProperties" -> false
)
val chat = new OpenAIChatCompletion().setResponseSchema(schema)
chat.setResponseSchema(schema, "answer_schema", strict = true)

OpenAIChatCompletion still returns the response envelope. Use Spark's built-in from_json on choices[0].message.content to create typed columns, or configure OpenAIPrompt with JSON post-processing. setResponseSchema does not change the output column's type.

ResponseFormat Options Quick Tests​

Below we validate text, json_object, and json_schema across APIs in separate cells for clarity.

from synapse.ml.services.openai import OpenAIChatCompletion
from pyspark.sql import Row

messages_df = spark.createDataFrame(
[([Row(role="user", content="Say hello")],)], ["messages"]
)
chat_text = (
OpenAIChatCompletion()
.setSubscriptionKey(key)
.setDeploymentName(deployment_name)
.setCustomServiceName(service_name)
.setMessagesCol("messages")
.setResponseFormat("text")
.setOutputCol("chat_text")
)
display(
chat_text.transform(messages_df)
.select("chat_text.choices.message.content")
.show(truncate=False)
)
messages_df = spark.createDataFrame(
[([Row(role="user", content="Return a JSON object with key greeting")],)],
["messages"],
)
chat_json_obj = (
OpenAIChatCompletion()
.setSubscriptionKey(key)
.setDeploymentName(deployment_name)
.setCustomServiceName(service_name)
.setMessagesCol("messages")
.setResponseFormat("json_object")
.setOutputCol("chat_json_obj")
)
display(
chat_json_obj.transform(messages_df)
.select("chat_json_obj.choices.message.content")
.show(truncate=False)
)

Responses API Structured JSON with json_schema (flattened form)​

For the Responses API, you can pass a flattened json_schema dict: top-level name, strict, and schema keys.

from synapse.ml.services.openai import OpenAIPrompt

schema = {
"type": "object",
"properties": {"answer": {"type": "string"}},
"required": ["answer"],
"additionalProperties": False,
}
df = spark.createDataFrame(
[("France", "capital"), ("Germany", "capital")], ["text", "category"]
)
prompt_flat = (
OpenAIPrompt()
.setSubscriptionKey(key)
.setDeploymentName(deployment_name)
.setCustomServiceName(service_name)
.setApiType("responses")
.setApiVersion("2025-04-01-preview")
.setPromptTemplate("What is the {category} of {text}.")
.setResponseFormat(
{
"type": "json_schema",
"name": "answer_schema",
"strict": True,
"schema": schema,
}
)
.setOutputCol("out_flat")
)
display(prompt_flat.transform(df).select("out_flat").show(truncate=False))

Invalid: Bare json_schema String​

Attempting to set responseFormat to the bare string json_schema should raise an error.

try:

OpenAIPrompt().setResponseFormat("json_schema")

except Exception as e:

print("Expected error:", e)

Working with Usage Statistics​

The following examples show how to enable usage tracking by setting the usageCol parameter. When set, usage statistics are returned in a separate column alongside the model output.

Chat Completions with usageCol​

Set the usageCol parameter to get usage statistics in a separate column.

from synapse.ml.services.openai import OpenAIPrompt

chat_usage_prompt = (
OpenAIPrompt()
.setSubscriptionKey(key)
.setDeploymentName(deployment_name)
.setCustomServiceName(service_name)
.setApiType("chat_completions")
.setUsageCol("usage") # Enable usage tracking with separate column
.setPromptTemplate("Provide a fun fact about {topic}.")
.setOutputCol("chat_output")
)

chat_usage_df = spark.createDataFrame([("Azure AI",)], ["topic"])
chat_usage_result = chat_usage_prompt.transform(chat_usage_df)

# Response is in chat_output, usage statistics are in the separate 'usage' column
display(chat_usage_result.select("chat_output", "usage").show(truncate=False))

Responses API with usageCol​

responses_usage_prompt = (
OpenAIPrompt()
.setSubscriptionKey(key)
.setDeploymentName(deployment_name)
.setCustomServiceName(service_name)
.setApiType("responses")
.setUsageCol("usage") # Enable usage tracking
.setPromptTemplate("List two key capabilities of {topic}.")
.setOutputCol("responses_output")
)

responses_usage_df = spark.createDataFrame([("Azure OpenAI",)], ["topic"])
responses_usage_result = responses_usage_prompt.transform(responses_usage_df)

# Response is in responses_output, usage is in separate 'usage' column
display(responses_usage_result.select("responses_output", "usage").show(truncate=False))

Embeddings with usageCol​

from synapse.ml.services.openai import OpenAIEmbedding

embedding_usage = (
OpenAIEmbedding()
.setSubscriptionKey(key)
.setCustomServiceName(service_name)
.setDeploymentName(deployment_name_embeddings)
.setUsageCol("usage") # Enable usage tracking with separate column
.setTextCol("text")
.setOutputCol("embeddings")
)

embedding_usage_df = spark.createDataFrame(
[("Usage statistics help monitor token consumption.",)], ["text"]
)
embedding_usage_result = embedding_usage.transform(embedding_usage_df)

# Embeddings are in 'embeddings' column, usage is in separate 'usage' column
display(embedding_usage_result.select("embeddings", "usage").show(truncate=False))

Responses API with store and Response Chaining​

The Responses API supports storing generated model responses for later retrieval via API using the store parameter. When store=True, the response ID is returned in a separate column (auto-generated or custom via setResponseIdCol). You can chain responses together using previousResponseId (scalar) or previousResponseIdCol (per-row from column) to maintain conversation context across requests.

from synapse.ml.services.openai import OpenAIResponses
from pyspark.sql.functions import col, lit, array, struct, expr
from pyspark.sql import Row

# Create messages DataFrame
messages_df = spark.createDataFrame(
[
([Row(role="user", content="What is Python?")],),
([Row(role="user", content="What is Spark?")],),
],
["messages"],
)

# First request with store=True to save response for later use
responses = (
OpenAIResponses()
.setSubscriptionKey(key)
.setDeploymentName(deployment_name)
.setCustomServiceName(service_name)
.setApiVersion("2025-04-01-preview")
.setStore(True)
.setMessagesCol("messages")
.setOutputCol("response")
)

first_result = responses.transform(messages_df).cache()
first_result.count() # Materialize to avoid re-execution

# Show response IDs (can be used for later retrieval)
display(
first_result.select(
col("response.id").alias("response_id"),
expr(
"element_at(filter(element_at(response.output, -1).content, x -> x.text is not null), 1).text"
).alias("text"),
).show(truncate=80)
)

# Chain responses using previousResponseIdCol
with_ids = first_result.withColumn("prev_id", col("response.id"))
follow_up_df = with_ids.select(
col("prev_id"),
array(
struct(lit("user").alias("role"), lit("Give me an example").alias("content"))
).alias("messages"),
)

chained = (
OpenAIResponses()
.setSubscriptionKey(key)
.setDeploymentName(deployment_name)
.setCustomServiceName(service_name)
.setApiVersion("2025-04-01-preview")
.setPreviousResponseIdCol("prev_id")
.setMessagesCol("messages")
.setOutputCol("chained_response")
)

chained_result = chained.transform(follow_up_df).cache()
display(
chained_result.select(
col("prev_id").alias("previous_response_id"),
expr(
"element_at(filter(element_at(chained_response.output, -1).content, x -> x.text is not null), 1).text"
).alias("chained_text"),
).show(truncate=80)
)

OpenAIPrompt with store and Response Chaining​

You can also use store with OpenAIPrompt for a simpler API. When store=True, the response ID is returned in a separate column (auto-generated by default or custom via setResponseIdCol). The main output column contains only the response text.

from synapse.ml.services.openai import OpenAIPrompt
from pyspark.sql.functions import col

# Sample DataFrame
prompt_df = spark.createDataFrame([("Python",), ("Spark",)], ["topic"])

# First request with store=True - response ID in separate auto-generated column
prompt_store = (
OpenAIPrompt()
.setSubscriptionKey(key)
.setDeploymentName(deployment_name)
.setCustomServiceName(service_name)
.setApiType("responses")
.setStore(True)
.setPromptTemplate("Briefly explain {topic}")
.setOutputCol("output")
)

first_prompt_result = prompt_store.transform(prompt_df).cache()
first_prompt_result.count()

# Get the auto-generated responseIdCol name
response_id_col_name = prompt_store.getResponseIdCol()

# Response text is in 'output', response ID is in auto-generated column
display(
first_prompt_result.select(
col("topic"),
col("output").alias("response"), # Just the text
col(response_id_col_name).alias("response_id"), # ID in separate column
).show(truncate=80)
)

# Chain using previousResponseIdCol
chained_prompt = (
OpenAIPrompt()
.setSubscriptionKey(key)
.setDeploymentName(deployment_name)
.setCustomServiceName(service_name)
.setApiType("responses")
.setPreviousResponseIdCol(response_id_col_name)
.setPromptTemplate("Give me a code example")
.setOutputCol("chained_output")
)

chained_prompt_result = chained_prompt.transform(first_prompt_result).cache()
display(
chained_prompt_result.select(
col("topic"),
col(response_id_col_name).alias("previous_id"),
col("chained_output").alias("follow_up_response"),
).show(truncate=80)
)